آموزش کامل Apache Kafka در ASP.NET Core
⚡ خلاصه سریع — اصل کاریها
اینها مهمترین چیزهایی هستن که اگه هیچی ندونی، فقط اینا رو بدونی کافیه:
- نکته ۱: Kafka یک Message Broker توزیعشده است که پیامها رو بهصورت جریانی (Stream) ذخیره و مدیریت میکنه.
- نکته ۲: دادهها در Topic ذخیره میشن و هر Topic میتونه به چند Partition تقسیم بشه.
- نکته ۳: Producer پیام تولید میکنه، Consumer پیام مصرف میکنه.
- نکته ۴: Consumer Group باعث میشه چند Consumer بار پردازش رو بین خودشون تقسیم کنن.
- نکته ۵: Offset شماره هر پیام تو Partition هست و Consumer با اون میدونه کجا مونده.
- نکته ۶: پیامها بعد از مصرف حذف نمیشن (تا زمان Retention)، این تفاوت اصلی Kafka با RabbitMQ است.
- نکته ۷: برای کار با Kafka در .NET از پکیج Confluent.Kafka استفاده میشه.
- نکته ۸: Kafka روی Zookeeper یا حالت جدید KRaft اجرا میشه.
- نکته ۹: Kafka داده رو روی دیسک ذخیره میکنه و بهخاطر Sequential I/O خیلی سریعه.
- نکته ۱۰: Kafka حداقل یهبار تحویل (At-least-once) تضمین میکنه، نه دقیقاً یهبار.
۳ اشتباه رایج که باید ازشون پرهیز کنی:
- Offset رو خودت Commit نمیکنی و پیامها گم میشن یا دوبار پردازش میشن.
- تعداد Partitionها رو بیدلیل زیاد میکنی که overhead ایجاد میکنه.
- خطاهای Deserialization رو تو Consumer هندل نمیکنی و کل برنامه هنگ میکنه.
اگر فقط ۵ دقیقه وقت داری، اینها رو یاد بگیر:
- Kafka = Topic + Producer + Consumer (سادهترین تصویر ذهنی)
dotnet add package Confluent.Kafka(پکیج اصلی کار با Kafka در .NET)- یه Producer بساز، پیام بفرست، یه Consumer بساز، پیام رو بخون
۱. مقدمه
این موضوع چیه؟
Apache Kafka یه پلتفرم توزیعشده برای Event Streaming است. به زبان سادهتر، یه صف پیام (Message Queue) خیلی قوی و سریعه که میتونه میلیونها پیام در ثانیه رو پردازش کنه.
تصور کن یه رستوران خیلی شلوغ داری:
- آشپز = Producer (پیام تولید میکنه — مثلاً یه سفارش غذا)
- گارسون = Kafka (پیام رو میگیره و تا زمان مصرف نگه میداره)
- میزبان = Consumer (پیام رو دریافت و پردازش میکنه)
تفاوت Kafka با صفهای معمولی اینه که گارسون (Kafka) غذا رو دور نمیریزه! اون نگهش میداره. اگه چندتا میزبان (Consumer) بیان، هرکدوم یه بخش از غذاها رو میگیرن و اگه یه میزبان غیبت کنه، غذاها هنوز هستن.
چرا باید یادش بگیری؟
| دلیل | توضیح |
|---|---|
| سرعت بالا | میتونه میلیونها پیام در ثانیه پردازش کنه |
| پایداری (Durability) | پیامها روی دیسک ذخیره میشن، حتی اگه سرور crash کنه |
| مقیاسپذیری | میتونی چندتا Broker اضافه کنی و بار رو تقسیم کنی |
| Replay | میتونی پیامهای قدیمی رو دوباره بخونی |
| تقاضای بازار | شرکتهای بزرگی مثل LinkedIn، Netflix، Uber ازش استفاده میکنن |
کجا استفاده میشه؟
- سیستمهای Microservices — ارتباط بین سرویسها بدون وابستگی مستقیم
- Log Aggregation — جمعآوری لاگ از سرورهای مختلف
- Event Sourcing — ذخیره تغییرات state بهصورت Event
- Real-time Analytics — تحلیل دادهها در لحظه
- Stream Processing — پردازش جریان دادهها
- نبود Single Point of Failure — اگه یه سرور بده، بقیه کار میکنن
۲. پیشنیازها
چه چیزهایی رو باید از قبل بدونی
| مهارت | سطح مورد نیاز |
|---|---|
| C# | متوسط (LINQ، async/await، Generics) |
| ASP.NET Core | مقدماتی (Minimal API یا Controller) |
| Docker | مقدماتی (run، compose) |
| Command Line | مقدماتی |
| مفاهیم Messaging | آشنایی کلی |
ابزارهای لازم
- .NET SDK 8.0 یا بالاتر — دانلود از سایت مایکروسافت
- Docker Desktop — برای اجرای Kafka بهصورت محلی
- Visual Studio 2022 یا VS Code یا JetBrains Rider
- kafka-topics CLI — ( داخل Docker هست، نیاز جداگانه نیست )
۳. نصب و راهاندازی
گام ۱: اجرای Kafka با Docker Compose
یه فایل به اسم docker-compose.yml بساز:
# docker-compose.yml
version: '3.8'
services:
kafka:
image: bitnami/kafka:latest
ports:
- "9092:9092"
environment:
# KRaft mode (بدون نیاز به Zookeeper)
KAFKA_CFG_NODE_ID: 0
KAFKA_CFG_PROCESS_ROLES: controller,broker
KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_CFG_LISTENERS_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
volumes:
- kafka_data:/bitnami/kafka
volumes:
kafka_data:
حالا اجراش کن:
docker-compose up -d
توضیح: این پیکربندی از حالت KRaft استفاده میکنه که دیگه نیازی به Zookeeper نداره. Kafka خودش مدیریت Controller رو انجام میده. خیلی سادهتر و سبکتره.
گام ۲: بررسی سلامت Kafka
docker exec -it <container_id> kafka-topics.sh --list --bootstrap-server localhost:9092
اگه بدون خطا اجرا شد و لیست خالی برگردوند، یعنی Kafka بالا اومده.
گام ۳: ساخت اولین پروژه ASP.NET Core
# ساخت solution
dotnet new sln -n KafkaTutorial
# ساخت Web API project
dotnet new webapi -n KafkaTutorial.Api -o src/KafkaTutorial.Api
# اضافه کردن پروژه به solution
dotnet sln add src/KafkaTutorial.Api
# نصب پکیج Kafka
cd src/KafkaTutorial.Api
dotnet add package Confluent.Kafka
توضیح: پکیج
Confluent.Kafkaکتابخانه رسمیه که Confluent (شرکت پشت Kafka) برای .NET منتشر میکنه. روی librdkafka (کتابخانه C) ساخته شده و performance بالایی داره.
۴. مفاهیم پایه
بیاید قبل از کد، مفاهیم رو با مثال ساده درک کنیم.
۴.۱ Broker
Broker یه سرور Kafka است. اگه ۳ تا سرور Kafka داشته باشی، ۳ تا Broker داری. Brokerها با هم صحبت میکنن و دادهها رو بین خودشون توزیع میکنن.
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Broker 1 │ │ Broker 2 │ │ Broker 3 │
│ :9092 │←→│ :9093 │←→│ :9094 │
└──────────┘ └──────────┘ └──────────┘
۴.۲ Topic
Topic یه دستهبندی پیامهاست. مثل یه کانال تلگرام — همه پیامهای مربوط به یه موضوع اونجا قرار میگیرن.
مثال: orders، user-events، payment-notifications
۴.۳ Partition
هر Topic به چندتا Partition تقسیم میشه. هر Partition یه log مرتب و تغییرناپذیر (Append-only) از پیامهاست.
Topic: orders
├── Partition 0: [msg1, msg2, msg5, msg8]
├── Partition 1: [msg3, msg4, msg6, msg9]
└── Partition 2: [msg7, msg10, msg11]
چرا Partition مهمه؟ چون Partition واحد موازیسازی است. اگه ۳ Partition داشته باشی، میتونی ۳ Consumer موازی داشته باشی که هرکدوم یه Partition رو میخونن.
۴.۴ Offset
هر پیام تو یه Partition یه شماره داره به اسم Offset. از ۰ شروع میشه و یکی یکی اضافه میشه.
Partition 0: [offset 0] [offset 1] [offset 2] [offset 3]
Consumer با ذخیره Offset فعلی، میدونه کجا مونده و بعداً از همونجا ادامه میده.
۴.۵ Producer
Producer پیام رو به Topic میفرسته. خودش تصمیم میگیره (یا با یه Key) پیام به کدوم Partition بره.
۴.۶ Consumer
Consumer پیامها رو از Topic میخونه. میتونه از یه Offset خاص شروع کنه.
۴.۷ Consumer Group
یه گروه از Consumerها که با هم کار میکنن و Partitionها رو بین خودشون تقسیم میکنن.
Consumer Group: "order-processors"
├── Consumer A ← reads Partition 0
├── Consumer B ← reads Partition 1
└── Consumer C ← reads Partition 2
قانون مهم: اگه تعداد Consumer بیشتر از Partition باشه، بعضی Consumerها هیچ کاری نمیکنن (idle میشن).
۴.۸ Retention
Kafka پیامها رو برای یه مدت مشخص نگه میداره (پیشفرض ۷ روز). بعد از اون، خودکار حذف میشن.
۴.۹ تجسم کلی
Producer Kafka Cluster Consumer
┌──────┐ ┌───────────────────────────┐ ┌──────────┐
│ App │────→│ Topic: orders │ ────────→│ Consumer │
│ │ │ ├── Partition 0 │ │ Group A │
└──────┘ │ ├── Partition 1 │ └──────────┘
│ └── Partition 2 │ ┌──────────┐
└───────────────────────────┘ ────────→│ Consumer │
│ Group B │
└──────────┘
۵. مثالهای کد ساده
۵.۱ ساخت Topic با کد
اول یه فولدر Kafka تو پروژه بساز و بعد یه کلاس KafkaConfig:
// Kafka/KafkaConfig.cs
namespace KafkaTutorial.Api.Kafka;
public static class KafkaConfig
{
public const string BootstrapServers = "localhost:9092";
public const string TopicName = "orders";
}
توضیح: این کلاس تنظیمات اصلی Kafka رو نگه میداره.
BootstrapServersآدرس سرور Kafka است. وقتی توسعه محلی میکنی،localhost:9092کافیه.
۵.۲ Producer ساده — ارسال یه پیام
یه کلاس KafkaProducerService بساز:
// Kafka/KafkaProducerService.cs
using Confluent.Kafka;
namespace KafkaTutorial.Api.Kafka;
public class KafkaProducerService
{
private readonly IProducer<Null, string> _producer;
private readonly string _topic = KafkaConfig.TopicName;
public KafkaProducerService()
{
var config = new ProducerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers
};
_producer = new ProducerBuilder<Null, string>(config).Build();
}
public async Task SendMessageAsync(string message)
{
var result = await _producer.ProduceAsync(
_topic,
new Message<Null, string> { Value = message }
);
Console.WriteLine(
$"✅ پیام ارسال شد — " +
$"Topic: {result.Topic}, " +
$"Partition: {result.Partition}, " +
$"Offset: {result.Offset}"
);
}
}
توضیح:
IProducer<Null, string>یعنی Key پیامnullاست و Value از نوعstring.ProduceAsyncپیام رو به Kafka میفرسته و منتظر میشه تا تأیید بشه.resultاطلاعات کامل پیام رو برمیگردونه — کدوم Partition، Offset چقدر شده.
۵.۳ استفاده از Producer تو Minimal API
// Program.cs
using KafkaTutorial.Api.Kafka;
var builder = WebApplication.CreateBuilder(args);
builder.Services.AddSingleton<KafkaProducerService>();
var app = builder.Build();
app.MapPost("/orders", async (KafkaProducerService producer, string order) =>
{
await producer.SendMessageAsync(order);
return Results.Ok(new { Message = "سفارش ثبت شد", Order = order });
});
app.Run();
توضیح: الان یه endpoint ساختی که پیام رو بهصورت POST میگیره و به Kafka میفرسته. تست:
curl -X POST "http://localhost:5000/orders?order=pizza-margherita"
۵.۴ Consumer ساده — خواندن پیام
یه کلاس KafkaConsumerService بساز که بهصورت Background Service اجرا بشه:
// Kafka/KafkaConsumerService.cs
using Confluent.Kafka;
namespace KafkaTutorial.Api.Kafka;
public class KafkaConsumerService : BackgroundService
{
private readonly string _topic = KafkaConfig.TopicName;
private readonly string _groupId = "order_processor_group";
private readonly string _bootstrapServers = KafkaConfig.BootstrapServers;
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
return Task.Run(() => StartConsumerLoop(stoppingToken), stoppingToken);
}
private void StartConsumerLoop(CancellationToken cancellationToken)
{
var config = new ConsumerConfig
{
BootstrapServers = _bootstrapServers,
GroupId = _groupId,
AutoOffsetReset = AutoOffsetReset.Earliest
};
using var consumer = new ConsumerBuilder<Null, string>(config).Build();
consumer.Subscribe(_topic);
try
{
while (!cancellationToken.IsCancellationRequested)
{
var consumeResult = consumer.Consume(cancellationToken);
Console.WriteLine(
$"📩 پیام دریافت شد — " +
$"Partition: {consumeResult.Partition}, " +
$"Offset: {consumeResult.Offset}, " +
$"Value: {consumeResult.Message.Value}"
);
}
}
catch (OperationCanceledException)
{
Console.WriteLine("Consumer متوقف شد.");
}
finally
{
consumer.Close();
}
}
}
توضیح:
BackgroundServiceیعتی این سرویس تو پسزمینه اجرا میشه.AutoOffsetReset.Earliestیعنی اولین باری که Consumer اجرا میشه، از ابتدای Topic بخونه.Consumeblocking است — منتظر میمونه تا پیام بیاد.consumer.Close()یه Commit خودکار از Offset فعلی انجام میده.
۵.۵ ثبت Consumer تو Program.cs
// Program.cs (بخشهای اضافهشده)
builder.Services.AddHostedService<KafkaConsumerService>();
۵.۶ تست کار
حالا برنامه رو اجرا کن:
dotnet run
یه ترمینال دیگه باز کن و پیام بفرست:
curl -X POST "http://localhost:5000/orders?order=pizza-margherita"
curl -X POST "http://localhost:5000/orders?order=burger"
curl -X POST "http://localhost:5000/orders?order=salad"
تو Console برنامه باید ببینی:
✅ پیام ارسال شد — Topic: orders, Partition: 0, Offset: 0
✅ پیام ارسال شد — Topic: orders, Partition: 0, Offset: 1
✅ پیام ارسال شد — Topic: orders, Partition: 0, Offset: 2
📩 پیام دریافت شد — Partition: 0, Offset: 0, Value: pizza-margherita
📩 پیام دریافت شد — Partition: 0, Offset: 1, Value: burger
📩 پیام دریافت شد — Partition: 0, Offset: 2, Value: salad
۶. مثال واقعی و کاربردی
پروژه کوچک: سیستم سفارشگذاری
سناریو: یه رستوران آنلاین داریم. وقتی مشتری سفارش میده:
OrderServiceیه پیام به TopicordersمیفرستهKitchenService(Consumer) سفارش رو دریافت و پردازش میکنهNotificationService(Consumer دیگه) ایمیل تأیید میفرسته
گام ۱: مدل داده
// Models/Order.cs
using System.Text.Json.Serialization;
namespace KafkaTutorial.Api.Models;
public class Order
{
public Guid Id { get; set; } = Guid.NewGuid();
public string CustomerName { get; set; } = string.Empty;
public string Item { get; set; } = string.Empty;
public decimal Amount { get; set; }
public DateTime CreatedAt { get; set; } = DateTime.UtcNow;
[JsonIgnore]
public string Status { get; set; } = "Pending";
}
گام ۲: Producer تایپشده با JSON Serialization
// Kafka/KafkaOrderProducer.cs
using System.Text.Json;
using Confluent.Kafka;
using KafkaTutorial.Api.Models;
namespace KafkaTutorial.Api.Kafka;
public class KafkaOrderProducer
{
private readonly IProducer<Null, string> _producer;
private readonly string _topic = KafkaConfig.TopicName;
public KafkaOrderProducer()
{
var config = new ProducerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers,
Acks = Acks.All, // منتظر تأیید همه Replicaها
EnableIdempotence = true // جلوگیری از ارسال دوباره
};
_producer = new ProducerBuilder<Null, string>(config).Build();
}
public async Task<(bool Success, string? Error)> SendOrderAsync(Order order)
{
try
{
var json = JsonSerializer.Serialize(order);
var result = await _producer.ProduceAsync(
_topic,
new Message<Null, string> { Value = json }
);
Console.WriteLine(
$"✅ سفارش ارسال شد — " +
$"Order ID: {order.Id}, " +
$"Partition: {result.Partition}, " +
$"Offset: {result.Offset}"
);
return (true, null);
}
catch (Exception ex)
{
return (false, ex.Message);
}
}
}
توضیح:
Acks.Allیعنی پیام حتماً روی چندتا سرور (Replica) ذخیره بشه قبل از تأیید.EnableIdempotence = trueیعنی اگه پیام بهخاطر مشکل شبکه دوباره ارسال بشه، Kafka فقط یه نسخه رو نگه میداره.
گام ۳: Consumer با پردازش Order
// Kafka/KafkaOrderConsumer.cs
using System.Text.Json;
using Confluent.Kafka;
using KafkaTutorial.Api.Models;
namespace KafkaTutorial.Api.Kafka;
public class KafkaOrderConsumer : BackgroundService
{
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<KafkaOrderConsumer> _logger;
public KafkaOrderConsumer(
IServiceProvider serviceProvider,
ILogger<KafkaOrderConsumer> logger)
{
_serviceProvider = serviceProvider;
_logger = logger;
}
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
return Task.Run(() => StartConsuming(stoppingToken), stoppingToken);
}
private void StartConsuming(CancellationToken cancellationToken)
{
var config = new ConsumerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers,
GroupId = "kitchen_service_group",
AutoOffsetReset = AutoOffsetReset.Earliest,
EnableAutoCommit = false // خودمون Commit میکنیم
};
using var consumer = new ConsumerBuilder<Null, string>(config).Build();
consumer.Subscribe(KafkaConfig.TopicName);
_logger.LogInformation("🍽️ Kitchen Service شروع شد...");
try
{
while (!cancellationToken.IsCancellationRequested)
{
var consumeResult = consumer.Consume(cancellationToken);
try
{
var order = JsonSerializer.Deserialize<Order>(consumeResult.Message.Value);
if (order != null)
{
_logger.LogInformation(
"🍳 آمادهسازی سفارش: {OrderItem} برای {Customer} — مبلغ: {Amount}",
order.Item,
order.CustomerName,
order.Amount
);
// شبیهسازی پردازش
Thread.Sleep(500);
// Manual commit بعد از پردازش موفق
consumer.Commit(consumeResult);
}
}
catch (Exception ex)
{
_logger.LogError(ex, "خطا در پردازش پیام: {Message}", consumeResult.Message.Value);
// اینجا میتونی پیام رو به یه Dead Letter Topic بفرستی
}
}
}
catch (OperationCanceledException)
{
_logger.LogInformation("Kitchen Service متوقف شد.");
}
finally
{
consumer.Close();
}
}
}
توضیح:
EnableAutoCommit = falseیعنی خودمون بعد از پردازش موفق Offset رو Commit میکنیم. این خیلی مهمه!- اگه برنامه Crash کنه وسط پردازش، چون Commit نکردیم، پیام دوباره از همون Offset خونده میشه.
consumer.Commit(consumeResult)فقط بعد از موفقیت کامل اجرا میشه.
گام ۴: Notification Consumer (گروه متفاوت)
// Kafka/KafkaNotificationConsumer.cs
using System.Text.Json;
using Confluent.Kafka;
using KafkaTutorial.Api.Models;
namespace KafkaTutorial.Api.Kafka;
public class KafkaNotificationConsumer : BackgroundService
{
private readonly ILogger<KafkaNotificationConsumer> _logger;
public KafkaNotificationConsumer(ILogger<KafkaNotificationConsumer> logger)
{
_logger = logger;
}
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
return Task.Run(() => StartConsuming(stoppingToken), stoppingToken);
}
private void StartConsuming(CancellationToken cancellationToken)
{
var config = new ConsumerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers,
GroupId = "notification_service_group", // گروه متفاوت!
AutoOffsetReset = AutoOffsetReset.Earliest
};
using var consumer = new ConsumerBuilder<Null, string>(config).Build();
consumer.Subscribe(KafkaConfig.TopicName);
_logger.LogInformation("📧 Notification Service شروع شد...");
try
{
while (!cancellationToken.IsCancellationRequested)
{
var consumeResult = consumer.Consume(cancellationToken);
try
{
var order = JsonSerializer.Deserialize<Order>(consumeResult.Message.Value);
_logger.LogInformation(
"📧 ایمیل تأیید ارسال شد به {Customer} برای سفارش: {Item}",
order?.CustomerName,
order?.Item
);
}
catch (Exception ex)
{
_logger.LogError(ex, "خطا در Notification");
}
}
}
catch (OperationCanceledException)
{
_logger.LogInformation("Notification Service متوقف شد.");
}
finally
{
consumer.Close();
}
}
}
نکته کلیدی: هر دو Consumer (Kitchen و Notification) از یه Topic میخونن، چون GroupId متفاوت دارن، هرکدوم تمام پیامها رو دریافت میکنن. این تفاوت اصلی با RabbitMQ است — چندتا Consumer مستقل میتونن همون پیام رو ببینن.
گام ۵: ثبت تو Program.cs
// Program.cs (نهایی)
using KafkaTutorial.Api.Kafka;
var builder = WebApplication.CreateBuilder(args);
builder.Services.AddSingleton<KafkaOrderProducer>();
builder.Services.AddHostedService<KafkaOrderConsumer>();
builder.Services.AddHostedService<KafkaNotificationConsumer>();
var app = builder.Build();
app.MapPost("/orders", async (KafkaOrderProducer producer, Order order) =>
{
var (success, error) = await producer.SendOrderAsync(order);
if (success)
return Results.Ok(new { Message = "سفارش ثبت شد", OrderId = order.Id });
return Results.StatusCode(500, new { Error = error });
});
app.Run();
گام ۶: تست
curl -X POST http://localhost:5000/orders \
-H "Content-Type: application/json" \
-d '{
"customerName": "علی رضایی",
"item": "پیتزا مارگاریتا",
"amount": 120000
}'
خروجی تو Console:
info: ✅ سفارش ارسال شد — Order ID: xxx, Partition: 0, Offset: 0
info: 🍽️ Kitchen Service شروع شد...
info: 📧 Notification Service شروع شد...
info: 🍳 آمادهسازی سفارش: پیتزا مارگاریتا برای علی رضایی — مبلغ: 120000
info: 📧 ایمیل تأیید ارسال شد به علی رضایی برای سفارش: پیتزا مارگاریتا
۷. مباحث پیشرفته
۷.۱ Partitioning با Key
وقتی یه Key به پیام میدی، Kafka با Hash اون Key تصمیم میگیره پیام به کدوم Partition بره. این تضمین میکنه همه پیامهای یه Key همیشه تو یه Partition باشن.
// Kafka/KafkaKeyedProducer.cs
using Confluent.Kafka;
namespace KafkaTutorial.Api.Kafka;
public class KafkaKeyedProducer
{
private readonly IProducer<string, string> _producer;
public KafkaKeyedProducer()
{
var config = new ProducerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers
};
_producer = new ProducerBuilder<string, string>(config).Build();
}
public async Task SendWithKeyAsync(string key, string value)
{
var result = await _producer.ProduceAsync(
KafkaConfig.TopicName,
new Message<string, string> { Key = key, Value = value }
);
Console.WriteLine(
$"Key: {key} → Partition: {result.Partition}, Offset: {result.Offset}"
);
}
}
توضیح: اگه ۳ تا Partition داشته باشی و Key
"user-123"رو بفرستی، Kafka با Hash اون Key یه Partition مشخص (مثلاً Partition 1) انتخاب میکنه. همه پیامهای"user-123"همیشه همون Partition 1 میرن. این برای ترتیببندی پیامهای یه کاربر حیاتی است.
۷.۲ Multiple Partitions و موازیسازی
برای ساخت Topic با چند Partition از CLI استفاده کن:
docker exec -it <container_id> kafka-topics.sh \
--create \
--bootstrap-server localhost:9092 \
--topic orders \
--partitions 3 \
--replication-factor 1
توضیح: الان ۳ Partition داری. اگه ۳ Consumer تو یه Group داشته باشی، هرکدوم یه Partition رو میخونن و موازی کار میکنن.
۷.۳ Dead Letter Queue (DLQ)
وقتی پردازش یه پیام خطا میده، بهجای از دست دادنش، به یه Topic دیگه به اسم DLQ بفرست.
// Kafka/DlqHandler.cs
using Confluent.Kafka;
namespace KafkaTutorial.Api.Kafka;
public class DlqHandler
{
private readonly IProducer<Null, string> _producer;
private const string DlqTopic = "orders.DLQ";
public DlqHandler()
{
var config = new ProducerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers
};
_producer = new ProducerBuilder<Null, string>(config).Build();
}
public async Task SendToDlqAsync(string originalMessage, string error)
{
var dlqMessage = $"{originalMessage} | ERROR: {error}";
await _producer.ProduceAsync(DlqTopic, new Message<Null, string> { Value = dlqMessage });
Console.WriteLine($"⚠️ پیام به DLQ فرستاده شد: {error}");
}
}
توضیح: DLQ یه الگوی استاندارد تو سیستمهای Messaging است. بهت اجازه میده پیامهای مشکلدار رو بعداً بررسی و پردازش کنی، بدون اینکه کل Pipeline متوقف بشه.
۷.۴ Schema Registry و Avro
برای serialize پیامها بهصورت بهینه و version-safe، از Avro و Schema Registry استفاده میکنن.
dotnet add package Confluent.SchemaRegistry
dotnet add package Confluent.SchemaRegistry.Serdes.Avro
// مثال سریع (نیاز به Schema Registry Server داره)
var config = new ProducerConfig
{
BootstrapServers = "localhost:9092"
};
var schemaRegistryConfig = new SchemaRegistryConfig
{
Url = "http://localhost:8081"
};
var avroSerializerConfig = new AvroSerializerConfig
{
AutoRegisterSchemas = true
};
توضیح: Avro یه فرمت باینری فشرده است. JSON خوانا ولی حجمش زیاده. Avro حجم کم و schema management قوی داره. برای پروژههای بزرگ، حتماً استفاده کن.
۷.۵ Exactly-Once Semantics
var config = new ProducerConfig
{
BootstrapServers = KafkaConfig.BootstrapServers,
EnableIdempotence = true,
Acks = Acks.All,
MaxInFlight = 5,
MessageSendMaxRetries = int.MaxValue
};
توضیح:
EnableIdempotenceیعنی Kafka تضمین میکنه پیام دقیقاً یهبار ثبت بشه، حتی اگه شبکه دوباره ارسال کنه. این فقط برای Producer هست — برای Consumer باید خودت Idempotent بنویسی (مثلاً با چک کردن ID پیام).
۷.۶ Batch Processing
بهجای ارسال یهبهیه، میتونی batchبندی کنی:
public async Task SendBatchAsync(IEnumerable<string> messages)
{
var tasks = messages.Select(msg =>
_producer.ProduceAsync(KafkaConfig.TopicName,
new Message<Null, string> { Value = msg })
);
await Task.WhenAll(tasks);
}
توضیح: Batch کردن میتونه throughput رو ۱۰ برابر کنه. ولی مواظب باش — اگه batch خیلی بزرگ باشه، latency زیاد میشه.
۷.۷ Kafka Streams (سادهسازی)
Kafka Streams برای پردازش جریانی استفاده میشه. تو .NET میتونی از Streamiz.Kafka.Net استفاده کنی:
dotnet add package Streamiz.Kafka.Net
// مثال ساده Stream Processing
var config = new StreamConfig<StringSerDes, StringSerDes>
{
ApplicationId = "order-stream-app",
BootstrapServers = "localhost:9092"
};
var stream = new KafkaStream(config);
stream.Start();
توضیح: Kafka Streams یه API برای پردازش real-time دادههاست. میتونی aggregation، filtering، joins و... انجام بدی. بهجای نوشتن Consumer دستی، از API سطح بالاتر استفاده میکنی.
۸. اشتباهات رایج
۸.۱ فعال کردن Auto Commit بیدقت
// ❌ اشتباه — Auto Commit فعاله
var config = new ConsumerConfig
{
EnableAutoCommit = true,
AutoCommitIntervalMs = 5000
};
// ✅ درست — Manual Commit
var config = new ConsumerConfig
{
EnableAutoCommit = false
};
// بعد از پردازش موفق:
consumer.Commit(consumeResult);
چرا؟ با Auto Commit، ممکنه Offset قبل از پردازش کامل Commit بشه. اگه Crash کنی، پیام از دست میره.
۸.۲ نداشتن تلاش مجدد (Retry)
وقتی پردازش خطا میده، پیام رو رها نکن:
// ❌ اشتباه
catch (Exception ex)
{
_logger.LogError(ex, "خطا");
// ادامه میدی و پیام گم میشه
}
// ✅ درست
catch (Exception ex)
{
_logger.LogError(ex, "خطا در پردازش، ارسال به DLQ");
await _dlqHandler.SendToDlqAsync(consumeResult.Message.Value, ex.Message);
consumer.Commit(consumeResult); // بعد از ارسال به DLQ
}
۸.۳ تعداد Partition زیاد بدون دلیل
❌ 100 Partition برای 10 پیام در دقیقه
✅ 3 Partition کافیه
چرا؟ هر Partition یه file روی دیسک و یه overhead مدیریتی داره. Partition اضافی منابع رو هدر میده.
۸.۴ Blocking در Consume
// ❌ اشتباه — پردازش طولانی تو حلقه
var result = consumer.Consume(cancellationToken);
Thread.Sleep(60000); // پردازش طولانی
consumer.Commit(result);
// ✅ درست — پردازش async
var result = consumer.Consume(cancellationToken);
await ProcessAsync(result.Message.Value);
consumer.Commit(result);
۸.۵ یه Partition با چند Consumer تو یه Group
اگه ۱ Partition داشته باشی و ۳ Consumer تو یه Group، ۲ تا Consumer idle میشن. این اشتباهه!
۸.۶ Ignoring Partitioner
وقتی به ترتیب پیامهای یه کاربر اهمیت داره، باید همیشه از Key استفاده کنی. بدون Key، ممکنه پیامهای یه کاربر تو Partitionهای مختلف پخش بشن و ترتیب بهم بخوره.
۹. بهترین روشها (Best Practices)
۹.۱ Use Dependency Injection
// ✅ اینجوری بنویس
builder.Services.AddSingleton<KafkaOrderProducer>();
builder.Services.AddHostedService<KafkaOrderConsumer>();
۹.۲ Use Configuration
// appsettings.json
{
"Kafka": {
"BootstrapServers": "localhost:9092",
"TopicName": "orders",
"ConsumerGroup": "order_processor_group"
}
}
// Kafka/KafkaSettings.cs
public class KafkaSettings
{
public string BootstrapServers { get; set; } = string.Empty;
public string TopicName { get; set; } = string.Empty;
public string ConsumerGroup { get; set; } = string.Empty;
}
// Program.cs
builder.Services.Configure<KafkaSettings>(builder.Configuration.GetSection("Kafka"));
۹.۳ همیشه Logging داشته باش
private readonly ILogger<KafkaOrderConsumer> _logger;
public async Task ProcessAsync(string message)
{
_logger.LogInformation("شروع پردازش پیام: {Message}", message);
try
{
// ...
_logger.LogInformation("پردازش موفق");
}
catch (Exception ex)
{
_logger.LogError(ex, "خطا در پردازش");
throw;
}
}
۹.۴ Monitoring و Metrics
از Confluent Control Center، Kafka UI، یا Prometheus + Grafana استفاده کن.
# سادهترین ابزار: Kafka UI
docker run -d -p 8080:8080 \
-e KAFKA_CLUSTERS_0_NAME=local \
-e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=host.docker.internal:9092 \
provectuslabs/kafka-ui:latest
باز کن: http://localhost:8080
۹.۵ Partition Count رو از اول درست انتخاب کن
| نرخ پیام | Partition پیشنهادی |
|---|---|
| < 100/sec | 1 |
| 100-1000/sec | 3-6 |
| 1000-10000/sec | 6-12 |
| > 10000/sec | 12+ |
نکته: بعد از ساخت Topic، اضافه کردن Partition امکانپذیره ولی ترتیب پیامها بهم میخوره. پس از همون اول درست تصمیم بگیر.
۹.۶ Idempotent Consumer بنویس
public async Task ProcessOrderAsync(Order order)
{
// چک کن آیا قبلاً پردازش شده؟
if (await _orderRepository.ExistsAsync(order.Id))
{
_logger.LogInformation("سفارش قبلاً پردازش شده: {OrderId}", order.Id);
return;
}
await _orderRepository.SaveAsync(order);
}
چرا؟ چون Kafka حداقل یهبار تحویل میده. ممکنه یه پیام دوبار بیاد. اگه Consumerت Idempotent نباشه، یه سفارش دو بار ثبت میشه!
۹.۷ Graceful Shutdown
public class KafkaOrderConsumer : BackgroundService
{
private IConsumer<Null, string>? _consumer;
protected override Task ExecuteAsync(CancellationToken stoppingToken)
{
// ...
}
public override Task StopAsync(CancellationToken cancellationToken)
{
_consumer?.Close(); // Commit آخرین Offset و بستن تمیز
return base.StopAsync(cancellationToken);
}
}
۹.۸ Avro/Protobuf برای Production
JSON خوبه برای شروع، ولی تو Production از Avro یا Protobuf استفاده کن.
۱۰. منابع و ادامه مسیر
منابع رسمی
دورهها و آموزشها
ابزارهای مفید
- Kafka UI — رابط گرافیکی برای مشاهده Topics، Consumers و...
- Confluent Control Center — پلتفرم کامل مانیتورینگ
- kafkacat (kcat) — CLI ابزار برای تست سریع
- Schema Registry — مدیریت schemaهای Avro
مسیر ادامه
- اول: این آموزش رو با یه پروژه کوچک تمرین کن — یه Producer، یه Consumer، یه Topic.
- بعد: Schema Registry و Avro رو یاد بگیر.
- بعد: Kafka Streams رو بررسی کن.
- بعد: Kafka Connect برای اتصال به دیتابیسها و سیستمهای دیگه.
- بعد: ksqlDB برای query کردن جریان دادهها بهصورت SQL.
- نهایتاً: Kafka Cluster Management و امنیت (SASL، SSL).
تمرین پیشنهادی
- یه پروژه ASP.NET Core با ۳ تا سرویس بساز:
OrderService(Producer)KitchenService(Consumer Group A)InventoryService(Consumer Group B)
- ۳ تا Partition بساز و موازیسازی رو تست کن.
- یه Consumer رو Stop کن و دوباره Start کن — باید از همون Offset ادامه بده.
- یه پیام خراب بفرست و DLQ رو امتحان کن.
- یه Dashboard ساده با یه API بزن که تعداد پیامهای دریافتی رو نشون بده.
به پایان آموزش رسیدیم! اگه تا اینجا با دقت همراه بودی، الان میتونی یه سیستم Kafka واقعی با ASP.NET Core پیادهسازی کنی. مهمترین چیز: شروع کن و تمرین کن. Kafka یه ابزار قدرتمنده، ولی پشتش یه منحنی یادگیری داره. قدم به قدم جلو بری، خیلی زود حرفهای میشی.